🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / takserver / src / commonMain / kotlin / org / meshtastic / core / takserver / TAKServer.kt
Displaying Raw • Download
core/takserver/src/commonMain/kotlin/org/meshtastic/core/takserver/TAKServer.kt copilot/create-implementation-plan (228d872f) Text, 8.00 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tf0883e@fileTb4b4b4:Te6edf3SuppressTb4b4b4(Ta5d6ff"Ta5d6ffTooGenericExceptionCaughtTa5d6ff"Tb4b4b4)
Tff7b72package T7ee787org.meshtastic.core.takserver
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787io.ktor.network.selector.SelectorManager
Tff7b72import T7ee787io.ktor.network.sockets.ServerSocket
Tff7b72import T7ee787io.ktor.network.sockets.Socket
Tff7b72import T7ee787io.ktor.network.sockets.SocketAddress
Tff7b72import T7ee787io.ktor.network.sockets.aSocket
Tff7b72import T7ee787kotlinx.coroutines.CancellationException
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.Job
Tff7b72import T7ee787kotlinx.coroutines.delay
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.StateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asStateFlow
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787org.meshtastic.core.di.CoroutineDispatchers
Tff7b72import T7ee787kotlin.random.Random
Tff7b72import T7ee787kotlinx.coroutines.isActive Tff7b72as Te6edf3coroutineIsActive
Tff7b72class T56d364TAKServerTb4b4b4(Tff7b72private Tff7b72val Te6edf3dispatchersTb4b4b4: Te6edf3CoroutineDispatchersTb4b4b4, Tff7b72private Tff7b72val Te6edf3portTb4b4b4: Tffa657Int Tff7b72= Te6edf3DEFAULT_TAK_PORTTb4b4b4) Tb4b4b4{
Tff7b72private Tff7b72var Te6edf3serverSocketTb4b4b4: Te6edf3ServerSocket? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3selectorManagerTb4b4b4: Te6edf3SelectorManager? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3running Tff7b72= Tff7b72false
Tff7b72private Tff7b72var Te6edf3serverScopeTb4b4b4: Te6edf3CoroutineScope? Tff7b72= Tff7b72null
Tff7b72private Tff7b72var Te6edf3acceptJobTb4b4b4: Te6edf3Job? Tff7b72= Tff7b72null
Tff7b72private Tff7b72val Te6edf3connectionsMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3connections Tff7b72= Te6edf3mutableMapOfTff7b72<Tffa657StringTb4b4b4, Te6edf3TAKClientConnectionTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_connectionCount Tff7b72= Te6edf3MutableStateFlowTb4b4b4(T79c0ff0Tb4b4b4)
Tff7b72val Te6edf3connectionCountTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72> Tff7b72= Te6edf3_connectionCountTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
Tff7b72var Te6edf3onMessageTb4b4b4: Tb4b4b4(Tb4b4b4(Te6edf3CoTMessageTb4b4b4) Tff7b72-Tff7b72> Tffa657UnitTb4b4b4)Tff7b72? Tff7b72= Tff7b72null
Tff7b72suspend Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4)Tb4b4b4: Te6edf3ResultTff7b72<Tffa657UnitTff7b72> Tb4b4b4{
T8b949e// Double-start guard: prevents SelectorManager / ServerSocket leaks
Tff7b72if Tb4b4b4(Te6edf3runningTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Server already running on port Tffd700$Te6edf3portTa5d6ff" Tb4b4b4}
Tff7b72return Te6edf3ResultTb4b4b4.Te6edf3successTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4}
Tff7b72return Tff7b72try Tb4b4b4{
Te6edf3serverScope Tff7b72= Te6edf3scope
T8b949e// Close any stale SelectorManager before creating a new one
Te6edf3selectorManagerTff7b72?.Te6edf3closeTb4b4b4(Tb4b4b4)
Te6edf3selectorManager Tff7b72= Te6edf3SelectorManagerTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3defaultTb4b4b4)
Te6edf3serverSocket Tff7b72= Te6edf3aSocketTb4b4b4(Te6edf3selectorManagerTff7b72!!Tb4b4b4)Tb4b4b4.Te6edf3tcpTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3bindTb4b4b4(Te6edf3hostname Tff7b72= Ta5d6ff"Ta5d6ff127.0.0.1Ta5d6ff"Tb4b4b4, Te6edf3port Tff7b72= Te6edf3portTb4b4b4)
Te6edf3running Tff7b72= Tff7b72true
Te6edf3acceptJob Tff7b72= Te6edf3scopeTb4b4b4.Te6edf3launchTb4b4b4(Te6edf3dispatchersTb4b4b4.Te6edf3ioTb4b4b4) Tb4b4b4{ Te6edf3acceptLoopTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3ResultTb4b4b4.Te6edf3successTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to bind TAK Server to 127.0.0.1:Tffd700$Te6edf3portTa5d6ff" Tb4b4b4}
Te6edf3ResultTb4b4b4.Te6edf3failureTb4b4b4(Te6edf3eTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72suspend Tff7b72fun Td2a8ffacceptLoopTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3scope Tff7b72= Te6edf3serverScope Tff7b72?: Tff7b72return
Tff7b72while Tb4b4b4(Te6edf3running Tff7b72&Tff7b72& Te6edf3scopeTb4b4b4.Te6edf3coroutineIsActiveTb4b4b4) Tb4b4b4{
Tff7b72try Tb4b4b4{
Tff7b72val Te6edf3clientSocket Tff7b72= Te6edf3serverSocketTff7b72?.Te6edf3acceptTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3clientSocket Tff7b72!Tff7b72= Tff7b72nullTb4b4b4) Tb4b4b4{
Te6edf3handleConnectionTb4b4b4(Te6edf3clientSocketTb4b4b4)
Tb4b4b4}
T8b949e// No delay on the success path — accept() is already suspending
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3CancellationExceptionTb4b4b4) Tb4b4b4{
Tff7b72throw Te6edf3e
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK server accept loop iteration failedTa5d6ff" Tb4b4b4}
T8b949e// Back-off only in the error path
Te6edf3delayTb4b4b4(Te6edf3TAK_ACCEPT_LOOP_DELAY_MSTb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffhandleConnectionTb4b4b4(Te6edf3clientSocketTb4b4b4: Te6edf3SocketTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3scope Tff7b72= Te6edf3serverScope Tff7b72?: Tff7b72return
Tff7b72val Te6edf3endpoint Tff7b72= Te6edf3clientSocketTb4b4b4.Te6edf3remoteAddressTb4b4b4.Te6edf3toStringTb4b4b4(Tb4b4b4)
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3clientSocketTb4b4b4.Te6edf3remoteAddressTb4b4b4.Te6edf3isLoopbackTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK server rejected non-loopback connection from Tffd700$Te6edf3endpointTa5d6ff" Tb4b4b4}
Te6edf3clientSocketTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
Tff7b72return
Tb4b4b4}
Tff7b72val Te6edf3connectionId Tff7b72= Te6edf3RandomTb4b4b4.Te6edf3nextIntTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3toStringTb4b4b4(Te6edf3TAK_HEX_RADIXTb4b4b4)
Tff7b72val Te6edf3clientInfo Tff7b72= Te6edf3TAKClientInfoTb4b4b4(Te6edf3id Tff7b72= Te6edf3connectionIdTb4b4b4, Te6edf3endpoint Tff7b72= Te6edf3endpointTb4b4b4)
Tff7b72val Te6edf3connection Tff7b72=
Te6edf3TAKClientConnectionTb4b4b4(
Te6edf3socket Tff7b72= Te6edf3clientSocketTb4b4b4,
Te6edf3clientInfo Tff7b72= Te6edf3clientInfoTb4b4b4,
Te6edf3onEvent Tff7b72= Tb4b4b4{ Te6edf3event Tff7b72-Tff7b72> Te6edf3handleConnectionEventTb4b4b4(Te6edf3connectionIdTb4b4b4, Te6edf3eventTb4b4b4) Tb4b4b4}Tb4b4b4,
Te6edf3scope Tff7b72= Te6edf3scopeTb4b4b4,
Tb4b4b4)
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3connectionsMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3connectionsTff7b72[Te6edf3connectionIdTff7b72] Tff7b72= Te6edf3connection
Te6edf3_connectionCountTb4b4b4.Te6edf3value Tff7b72= Te6edf3connectionsTb4b4b4.Te6edf3size
Tb4b4b4}
Te6edf3connectionTb4b4b4.Te6edf3startTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tff7b72private Tff7b72fun Td2a8ffhandleConnectionEventTb4b4b4(Te6edf3connectionIdTb4b4b4: Tffa657StringTb4b4b4, Te6edf3eventTb4b4b4: Te6edf3TAKConnectionEventTb4b4b4) Tb4b4b4{
Tff7b72when Tb4b4b4(Te6edf3eventTb4b4b4) Tb4b4b4{
Tff7b72is Te6edf3TAKConnectionEventTb4b4b4.Te6edf3Message Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3onMessageTff7b72?.Te6edf3invokeTb4b4b4(Te6edf3eventTb4b4b4.Te6edf3cotMessageTb4b4b4)
Tb4b4b4}
Tff7b72is Te6edf3TAKConnectionEventTb4b4b4.Te6edf3Disconnected Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3serverScopeTff7b72?.Te6edf3launch Tb4b4b4{
Te6edf3connectionsMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3connectionsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3connectionIdTb4b4b4)
Te6edf3_connectionCountTb4b4b4.Te6edf3value Tff7b72= Te6edf3connectionsTb4b4b4.Te6edf3size
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72is Te6edf3TAKConnectionEventTb4b4b4.Te6edf3Error Tff7b72-Tff7b72> Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eventTb4b4b4.Te6edf3errorTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK client connection error: Tffd700$Te6edf3connectionIdTa5d6ff" Tb4b4b4}
Te6edf3serverScopeTff7b72?.Te6edf3launch Tb4b4b4{
Te6edf3connectionsMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Te6edf3connectionsTb4b4b4.Te6edf3removeTb4b4b4(Te6edf3connectionIdTb4b4b4)
Te6edf3_connectionCountTb4b4b4.Te6edf3value Tff7b72= Te6edf3connectionsTb4b4b4.Te6edf3size
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72is Te6edf3TAKConnectionEventTb4b4b4.Te6edf3Connected Tff7b72-Tff7b72> Tb4b4b4{
T8b949e/* no-op: logged by TAKClientConnection.start() */
Tb4b4b4}
Tff7b72is Te6edf3TAKConnectionEventTb4b4b4.Te6edf3ClientInfoUpdated Tff7b72-Tff7b72> Tb4b4b4{
T8b949e/* no-op: TAKClientConnection tracks updated info locally */
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4) Tb4b4b4{
Te6edf3running Tff7b72= Tff7b72false
Te6edf3acceptJobTff7b72?.Te6edf3cancelTb4b4b4(Tb4b4b4)
Te6edf3acceptJob Tff7b72= Tff7b72null
T8b949e// Close connections synchronously — TAKClientConnection.close() is non-suspending,
T8b949e// so we don't need to launch into the (possibly-cancelled) serverScope.
Tff7b72val Te6edf3toCloseTb4b4b4: Te6edf3ListTff7b72<Te6edf3TAKClientConnectionTff7b72>
T8b949e// We can't use Mutex.withLock here (non-suspending context) so we swap & clear under a
T8b949e// best-effort copy — worst case a connection added concurrently is closed by socket teardown.
Te6edf3toClose Tff7b72= Te6edf3connectionsTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4)
Te6edf3connectionsTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3_connectionCountTb4b4b4.Te6edf3value Tff7b72= T79c0ff0
Te6edf3toCloseTb4b4b4.Te6edf3forEach Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3serverSocketTff7b72?.Te6edf3closeTb4b4b4(Tb4b4b4)
Te6edf3serverSocket Tff7b72= Tff7b72null
Te6edf3selectorManagerTff7b72?.Te6edf3closeTb4b4b4(Tb4b4b4)
Te6edf3selectorManager Tff7b72= Tff7b72null
Te6edf3serverScope Tff7b72= Tff7b72null
Tb4b4b4}
Tff7b72suspend Tff7b72fun Td2a8ffbroadcastTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72val Te6edf3currentConnections Tff7b72= Te6edf3connectionsMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3connectionsTb4b4b4.Te6edf3valuesTb4b4b4.Te6edf3toListTb4b4b4(Tb4b4b4) Tb4b4b4}
Te6edf3currentConnectionsTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3connection Tff7b72-Tff7b72>
Tff7b72try Tb4b4b4{
Te6edf3connectionTb4b4b4.Te6edf3sendTb4b4b4(Te6edf3cotMessageTb4b4b4)
Tb4b4b4} Tff7b72catch Tb4b4b4(Te6edf3eTb4b4b4: Te6edf3ExceptionTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3wTb4b4b4(Te6edf3eTb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to broadcast CoT to TAK client Tffd700${Te6edf3connectionTb4b4b4.Te6edf3clientInfoTb4b4b4.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4}
Te6edf3connectionTb4b4b4.Te6edf3closeTb4b4b4(Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72suspend Tff7b72fun Td2a8ffhasConnectionsTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tff7b72= Te6edf3connectionsMutexTb4b4b4.Te6edf3withLock Tb4b4b4{ Te6edf3connectionsTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Returns true if this [SocketAddress] represents a loopback address (IPv4 127.x.x.x or IPv6 ::1).
*
* Ktor's [SocketAddress.toString] returns strings like "/127.0.0.1:4242" (JVM) or "127.0.0.1:4242" on other platforms,
* so we strip any leading slash and check prefixes without parsing the host. This keeps the check in commonMain without
* an expect/actual.
*/
Tff7b72private Tff7b72fun Te6edf3SocketAddressTb4b4b4.Td2a8ffisLoopbackTb4b4b4(Tb4b4b4)Tb4b4b4: Tffa657Boolean Tb4b4b4{
Tff7b72val Te6edf3addr Tff7b72= Te6edf3toStringTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3removePrefixTb4b4b4(Ta5d6ff"Ta5d6ff/Ta5d6ff"Tb4b4b4)
Tff7b72return Te6edf3addrTb4b4b4.Te6edf3startsWithTb4b4b4(Ta5d6ff"Ta5d6ff127.Ta5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Te6edf3addrTb4b4b4.Te6edf3startsWithTb4b4b4(Ta5d6ff"Ta5d6ff::1Ta5d6ff"Tb4b4b4) Tff7b72|Tff7b72| Te6edf3addrTb4b4b4.Te6edf3startsWithTb4b4b4(Ta5d6ff"Ta5d6ff[::1]Ta5d6ff"Tb4b4b4)
Tb4b4b4}
Served by rngit 1.5.1 - Generated in 0.07s